/** * Search-box.ai * Author: Simon-Pierre Boucher * Contact: contact@spboucher.ai * File: apps/web/app/api/research/[id]/stream/route.ts * Description: SSE stream of research events with Last-Event-ID replay (proxy-safe headers). */ import { events, sessions } from "@search-box/db"; import { ensureMigrated } from "@/lib/runner"; export const runtime = "nodejs"; export const dynamic = "force-dynamic"; const POLL_MS = 400; const HEARTBEAT_MS = 15_000; export async function GET( req: Request, { params }: { params: Promise<{ id: string }> } ): Promise { await ensureMigrated(); const { id } = await params; const session = await sessions.get(id); if (!session) return new Response("not found", { status: 404 }); // Replay position: Last-Event-ID header (EventSource reconnect) or ?from= const url = new URL(req.url); const fromParam = url.searchParams.get("from"); const lastEventId = req.headers.get("last-event-id"); let cursor = Number(lastEventId ?? fromParam ?? 0); if (!Number.isFinite(cursor) || cursor < 0) cursor = 0; const encoder = new TextEncoder(); const stream = new ReadableStream({ async start(controller) { let closed = false; const close = () => { if (!closed) { closed = true; try { controller.close(); } catch { /* already closed */ } } }; req.signal.addEventListener("abort", close); const send = (seq: number, data: string) => { controller.enqueue(encoder.encode(`id: ${seq}\ndata: ${data}\n\n`)); }; let lastBeat = Date.now(); try { while (!closed) { const batch = await events.listAfter(id, cursor); for (const ev of batch) { cursor = ev.seq; send(ev.seq, JSON.stringify(ev.payload)); } // Stop once the session reached a terminal state and all events were sent. if (batch.length === 0) { const current = await sessions.get(id); if (current && (current.status === "completed" || current.status === "failed")) { controller.enqueue(encoder.encode(`event: end\ndata: {}\n\n`)); break; } } if (Date.now() - lastBeat >= HEARTBEAT_MS) { controller.enqueue(encoder.encode(`: heartbeat\n\n`)); lastBeat = Date.now(); } await new Promise((r) => setTimeout(r, POLL_MS)); } } catch { // client disconnected or db error — just end the stream } finally { close(); } } }); return new Response(stream, { headers: { "Content-Type": "text/event-stream; charset=utf-8", "Cache-Control": "no-cache, no-transform", Connection: "keep-alive", "X-Accel-Buffering": "no" } }); }